Skip to content

SS-455 - source: commit Kafka offsets on change and on a timer - #39462

Merged
martykulma merged 2 commits into
maz-offset-committed-per-exportfrom
maz-kafka-offset-commit-dedupe
Oct 5, 2026
Merged

martykulma merged 2 commits into
maz-offset-committed-per-exportfrom
maz-kafka-offset-commit-dedupe

Conversation

@martykulma

@martykulma martykulma commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

Each worker commits only when its offsets change, plus on a timer. The timer is needed because the consumer only ever assigns partitions, which results in the brokers treating the group as standalone. For standalone groups, brokers expire a partition's offset offsets.retention.minutes after its last commit, whether or not it is still current. The interval is the kafka_offset_commit_refresh_interval dyncfg: 10 minutes by default, 10 seconds in CI, zero disables it. A failed commit is retried on the next resume upper.

🤖 Generated with Claude Code
👨 Improved by human

@martykulma
martykulma added this pull request to stack #39464 October 1, 2026 18:35
@martykulma martykulma changed the title source: Kafka only comits offsets for a partition if they've changed SS-455 - source: Kafka only commits offsets for a partition if they've changed Oct 1, 2026
@linear-code

linear-code Bot commented Oct 1, 2026

Copy link
Copy Markdown

SS-455

@martykulma
martykulma marked this pull request as ready for review October 1, 2026 19:30
@martykulma
martykulma requested a review from a team as a code owner October 1, 2026 19:30
@def-

def- commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

QA LLM Review

1. MEDIUM -- Idle partitions are never re-committed, so the broker expires their consumer-group offsets

src/storage/src/source/kafka.rs:1163

With the offsets != self.committed_offsets check, a worker whose partitions receive no new messages commits them once and never again, even while the rest of the source keeps advancing. Materialize consumes with assign(), which makes it a standalone consumer group, and brokers delete a standalone group's offset for a partition offsets.retention.minutes (7 days by default) after that partition's last commit. Lag-monitoring tools, which the Kafka source docs say these commits exist for, then lose those partitions.

Details

On main, reclock_committed_upper calls tx.send_replace every time it runs (source_reader_pipeline.rs:685 at the merge base), so every worker re-commits its offsets about once per tick, and that refreshes the commit timestamp. On #39461, any change to ResumeUppers still re-commits every worker's partitions. This PR takes away that last refresh for idle partitions. Two examples: a keyed topic where some partitions get no writes for a week, or a cluster with at least as many workers as partitions, where each worker owns at most one partition. After the retention window, kafka-consumer-groups --describe shows - for those partitions, and lag exporters either drop them or report them as missing. Data ingestion is not affected.

Suggested fix: keep the dedup, but re-commit unchanged offsets once the last successful commit is older than a refresh interval set well below typical broker retention (for example one hour, or a dyncfg). Because #39461 switched to send_if_modified, process_frontier is never called for a fully idle source. So the refresh needs a timer in resume_uppers_process_loop, for example a select! over the stream and a tokio::time::interval, rather than only a looser condition here. The same timer would also retry a failed commit on an idle source. Today the doc comment's "retried on the next resume upper" may never happen for such a source.

@martykulma
martykulma force-pushed the maz-kafka-offset-commit-dedupe branch from 7e97474 to e1893ed Compare October 2, 2026 14:21
@martykulma
martykulma requested a review from a team as a code owner October 2, 2026 14:21
@martykulma martykulma changed the title SS-455 - source: Kafka only commits offsets for a partition if they've changed SS-455 - source: commit Kafka offsets on change and on a timer Oct 2, 2026
@martykulma

Copy link
Copy Markdown
Contributor Author
  1. MEDIUM -- Idle partitions are never re-committed, so the broker expires their consumer-group offsets

Confirmed and addressed

@def-

def- commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

QA LLM Review

1. MEDIUM -- "10m" is not a valid system-parameter duration, so environmentd fails to boot in random-parameter CI runs

misc/python/materialize/mzcompose/__init__.py:428

System variables parse durations with the units us, ms, s, min, h and d. There is no m, so "10m" is rejected. Under CI_SYSTEM_PARAMETERS=random, about one job in three picks "10m", and every Materialize container in that job then panics at startup. All tests in those jobs fail, which hides whatever the random-parameter run was meant to find.

Details

Dyncfgs that become SQL system vars are parsed by impl Value for Duration (src/sql/src/session/vars/value.rs:307), not by humantime, and that parser has no m unit. get_default_system_parameters passes the chosen value as a system-parameter default. When catalog open hits a parse error it returns that error instead of skipping the parameter (src/adapter/src/catalog/open.rs:263).

I reproduced this on a current PR image using an existing duration dyncfg:

MZ_SYSTEM_PARAMETER_DEFAULT="kafka_poll_max_wait=10m"
environmentd: fatal: parameter "kafka_poll_max_wait" cannot have value "m": expected us, ms, s, min, h, or d but got "m"

Every other duration in this file already uses min (for example "1min" at line 578). Changing the value to "10min" fixes it.

@patrickwwbutler patrickwwbutler left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Small question, plus suggestion to fix the LLM lint, but no blockers! LGTM

# Low default so CI exercises the periodic recommit, which production
# only reaches after ten minutes.
VariableSystemParameter(
"kafka_offset_commit_refresh_interval", "10s", ["1s", "10s", "10m"]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
"kafka_offset_commit_refresh_interval", "10s", ["1s", "10s", "10m"]
"kafka_offset_commit_refresh_interval", "10s", ["1s", "10s", "600s"]

// Zero disables the refresh. `interval` panics on a zero period, so the tick arm
// below is what honors it.
let mut refresh =
tokio::time::interval(refresh_interval.max(Duration::from_millis(1)));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

if the value was set to 0, and this clamps it to 1ms, is that sustainable? I suppose it would just happen every iteration of the loop, which is okay?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1ms would just update as fast as the calls would allow it (because the missed tick behavior is delay).

If someone sets it to 0 to disable, then the branch never fires because the precondition would return false:

_ = refresh.tick(), if !refresh_interval.is_zero() => {
if let Err(e) = offset_committer.refresh().await {
report(e, &"refresh");
}

https://docs.rs/tokio/latest/tokio/macro.select.html

Evaluate all provided expressions. If the precondition returns false, disable the branch for the remainder of the current call to select!. Re-entering select! due to a loop clears the “disabled” state.

@martykulma
martykulma force-pushed the maz-kafka-offset-commit-dedupe branch from e1893ed to 97b4ca9 Compare October 5, 2026 17:56
martykulma and others added 2 commits October 5, 2026 14:13
Brokers expire a standalone group's offset per partition
offsets.retention.minutes after its last commit, so offsets the dedupe
skips must still be recommitted periodically. The interval is the
kafka_offset_commit_refresh_interval dyncfg, 10 minutes by default and
10 seconds in CI, with zero disabling the refresh.

Co-Authored-By: Claude <noreply@anthropic.com>
@martykulma
martykulma force-pushed the maz-kafka-offset-commit-dedupe branch from 97b4ca9 to c220506 Compare October 5, 2026 18:13
@martykulma
martykulma merged commit 491235b into main Oct 5, 2026
84 checks passed
@martykulma
martykulma deleted the maz-kafka-offset-commit-dedupe branch October 5, 2026 18:44
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants